iT邦幫忙

2026 iThome 鐵人賽

DAY 18
0

D16 講過 DataFusion 怎麼執行(拉式、batch、tokio)、D17 講記憶體資源的協調。今天要討論是什麼東西這麼關鍵,幾乎決定整條 SQL 查詢的成本與資源用量、planner 有最多路徑可以選?

JOIN

同一句 A JOIN B ON A.x = B.x,DataFusion 至少有四種完全不同的執行策略可以走:hash、sort-merge、symmetric hash、nested loop。四個的記憶體用量、磁碟寫入、CPU 需求、能不能串流全都不一樣。

Planner 憑什麼選?表面上是一棵很短的決策樹,底下是一整套統計 + 成本模型 + 動態調整。今天把這條路徑拆開,順便看一個很優雅的技巧 dynamic filter:規劃期看的資訊不夠,執行期把 build side 建好的東西「橫向」傳過去,讓 probe side 再剪一次。

一句 SQL 有幾種跑法

先看一個很日常的 join:

SELECT o.o_orderkey, c.c_name
FROM orders o JOIN customer c ON o.o_custkey = c.c_custkey
WHERE c.c_region = 'ASIA'

Planner 看到這句,要決定的第一件事不是「怎麼跑」,而是「這個 join 適用哪些演算法」。
DataFusion 主要選項有四個:

  • HashJoinExec:一邊建 hash table、另一邊 probe
  • SortMergeJoinExec:兩邊排好序,一起走
  • SymmetricHashJoinExec:兩邊都建 hash table,兩邊都當 probe
  • NestedLoopJoinExec:兩層迴圈掃到底

四個都能跑同一句 SQL,但成本差好幾個數量級,挑錯的代價不是「慢一點」,是「這條 query 跑不完」。

決策樹的第一分岔:有沒有等值謂詞

四個演算法裡有三個都要靠匹配 key:hash 打 bucket、sort-merge 兩邊一起走。
這些只能回答一種問題「這兩個值相等嗎」。

如果 JOIN 條件是 ><BETWEEN 這種範圍,或更誇張的 ON true(笛卡兒積),沒有 key 可以打,也沒有排序能一起走。這時候只剩下 nested loop。

所以第一問一定是:有沒有等值謂詞?

  • 沒有 → NestedLoopJoin(沒選擇)
  • 有 → 進第二問

Nested loop 的複雜度是 O(m × n),兩邊都是百萬列的時候直接不用想。實務上有個常見加速:broadcast NLJ,小邊 broadcast 到每個 partition、大邊在 partition 內走 loop。理論複雜度沒變,但小邊留在 CPU cache 內、常數項會小很多。

第二分岔:兩邊排序好了嗎

SortMergeJoin 需要兩邊輸入都已排序,乍看很吃虧為了 join 額外做 sort,比直接 hash 慢。
但它有兩個場景會贏:

  1. 輸入天然有序:資料本來就從有序 source 進來(bucketed table、上游剛做完 sort、Iceberg sorted table),sort 那步免了
  2. 兩邊都很大:hash 兩邊都塞不進記憶體、就算 spill 也是兩邊 spill;SortMerge 只要 sort 一次然後 merge,記憶體用量小很多

Spark 3.x 之後大 join 常常預設走 SortMerge 就是這個道理,事實表對事實表,兩邊都是幾億筆,hash 直接跪。

Comet 也 offload SMJ,跟 Spark 預設一致。

第三分岔:哪一邊可以先物化

如果沒有天然排序,就走 hash 這條路。但 hash 有個隱含假設:build 側可以先物化完
下游要 await build 建完才能開始 probe。

大部分批次查詢都成立,可是有些場景不行:

  • 串流查詢:兩邊都是 unbounded stream,等對面「物化完」永遠不會發生
  • 某些 outer join:full outer 要兩邊完整看過才知道 unmatched,單邊 hash 不夠

這時候用 SymmetricHashJoin,兩邊各建一份 hash table,一邊來一 batch 就 probe 另一邊 + 插入自己這邊。代價是兩份 hash table 的記憶體,比一般 hash 貴一倍以上。

DataFusion 主要拿這個處理串流,Comet 是批次引擎,這個一般不 offload。

第四分岔:hash 用哪一邊當 build

走到 hash 這條路,還剩最後一問哪一邊當 build、哪一邊當 probe

答案很直白:挑估算 row count 小的當 build。build side 要整個塞進 hash table,越小記憶體越省,probe side 越大 CPU 就一次消化越多。

但這裡藏了一個很重的地雷如果估算 build side 的 row count 錯誤,會導致記憶體不足、hash table 需要 spill 到 disk,效能嚴重下降,查詢時間可能多出數倍

如果統計不準、把大表當 build,hash table 就爆到要 spill。
回想 D17 中,HashJoin 的 build 側是 spillable 的,MemoryPool 拒絕 reservation 時它會把部分 hash bucket 寫到 disk。
Spill 之後 probe 進來每一筆都要先看記憶體、再看 disk,兩趟 I/O 疊上去。

所以你會看到 CBO(cost-based optimizer)這個詞在 join 場景特別重要——選錯 build side 的代價是幾倍,不是幾成。

決策樹全貌

把上面拼起來:

       ┌──────────────────────────┐
       │ 有等值謂詞嗎?             │
       └────┬──────────────┬──────┘
            │沒有          │有
            ▼              ▼
     NestedLoopJoin    兩邊都排序好?
     (或 broadcast)   ┌────┴────┐
                     │是       │否
                     ▼         ▼
              SortMergeJoin  兩邊都不能先物化?
                            ┌───┴───┐
                            │是     │否
                            ▼       ▼
                   SymmetricHashJoin  HashJoin
                                     (挑小的當 build)

實務上還有更多細節:left / right / full outer 各有偏好、broadcast vs shuffle 是另一個維度、null-aware anti join 有特殊處理。但主幹就這條。

規劃期的決策不夠:Dynamic Filter

上面決策樹都是 planner 在規劃期做的事,看 SQL、看統計、選演算法
可是有個問題:很多資訊要跑起來才知道

回到開頭那個例子:orders JOIN customer ON o.o_custkey = c.c_custkey WHERE c.c_region = 'ASIA'

規劃期只知道「customer 這邊會被 filter」,但 filter 完剩多少 key、剩下哪些 key,都得跑起來才知道。假設 c_region = 'ASIA' 篩完只剩 200 個 c_custkey——那 orders 本來要掃 6 億筆,其實只要看那 200 個 c_custkey 出現在哪些 row group 裡就好。

Dynamic Filter 就是這件事:

  1. HashJoin build 側跑完(拿到 customer 的 200 個 key)
  2. 把這 200 個 key 打包成一個 filter(bloom filter、min-max range、in-list 都可以)
  3. 動態下推給 probe 側的 scan
  4. Scan 用這個 filter 做 row group / page 剪枝

這個技巧的正式名字是 sideways information passing——資訊不是走 pipeline 上下游(那是 top-down),而是橫向從一個 operator 傳到另一個。

為什麼 dynamic filter 有效

為什麼把 join filter 動態下推到 scan 是有效的

粗糙版答案:「因為 scan 少讀了」但這個答案沒觸到重點,如果 scan 只是把資料讀進來後在 filter 節點過濾,那 dynamic filter 只是省了一段 in-memory 處理,不會有數量級差別。

真正有效的原因是 Parquet 的檔案內剪枝機制

  • Parquet footer 有 row group 層的 min-max 統計,dynamic filter 是 in-list 時可以直接跳過整個 row group
  • Parquet 1.12 起有 bloom filter,可以在 page 層剪枝
  • 剪掉的 row group / page 完全不 decode,I/O 也不做

也就是說 dynamic filter 有效不是因為它是個 filter,而是因為它讓 scan 少做 I/O 跟 decode。如果 probe side 是 CSV(沒有檔案內剪枝機制),這個技巧退化成只是 in-memory filter,效果差很多。

還有一個時機關鍵:dynamic filter 必須在 probe scan 還沒開始(或早期)時送到。太晚 scan 已經讀完就沒意義。DataFusion 用 async channel + repartition 邊界處理這件事,也是 pull + async 執行模型的一個副作用——channel 天然可以做橫向資訊傳遞。

回到 Comet:哪些 join 有 offload

Comet 目前的 join 支援大致是(細節看 release note 跟 compatibility matrix):

  • HashJoinExec(含 inner / left / right / left semi / left anti)—— offload
  • BroadcastHashJoinExec —— offload(維度表 join 很常見)
  • SortMergeJoinExec —— offload(Spark 大 join 的預設)
  • BroadcastNestedLoopJoinExec —— 部分 offload,要看 join type
  • SymmetricHashJoinExec —— 一般不 offload(串流不是 Comet 主打)
  • 純 NestedLoopJoinExec —— 一般 fallback

Fallback 的代價回到 D15 接點五——整個 join 子樹回到 JVM 走 Spark 原生,這一批要從 Arrow C Data Interface 轉回 UnsafeRow(付一次 C2R)、join 完再考慮要不要轉回 Arrow 繼續 native(付一次 R2C)。所以 fallback 不是「這個 join 慢一點」,是這一批資料要多兩次跨界

結論

Join 演算法選擇的主幹是短短一棵決策樹:有等值 → 有序 → 能物化 → 挑小邊。挑錯的代價不是慢,是查詢跑不完。Dynamic Filter 是規劃期的決策不夠時、執行期補一刀的技巧,有效性靠 Parquet 檔案內剪枝背書——是儲存跟執行的一次成功協作。

明天 D19 接執行細節的下一題:資料超過記憶體時排序怎麼繼續?外部排序是所有 pipeline breaker 都會遇到的降級路徑,順便把「下推」跟「提前終止」這幾個省 I/O 的手法一起說完。

那就明天見~

參考資料


上一篇
Day 17 MemoryPool:兩種公平
下一篇
Day 19 SortExec 記憶體不夠時怎麼寫磁碟:排序、Pushdown與 Top-K
系列文
1+1+1>3 ~ Spark 與 DataFusion、Comet 效能煉金術 ~21
圖片
  熱門推薦
圖片
{{ item.channelVendor }} | {{ item.webinarstarted }} |
{{ formatDate(item.duration) }}
直播中

尚未有邦友留言

立即登入留言